Repository navigation
Conversation
## Why Model authors should declare their worker and typed video output without reimplementing threading, session bookkeeping, and cleanup around WorkerGroup. This layer supplies that managed authoring path without introducing an engine-host dependency. ## What Changed Adds the experimental DistributedVideoModel adapter with worker configuration, control snapshots, frame-to-output mapping, pause/restart, and worker cleanup. Authoring and configuration tests cover the adapter contract, while the dependent example PR exercises the complete run-loop lifecycle. The scope remains single-node, single-session uint8 video, not arbitrary structured outputs or multi-session scheduling. Local Python 3.12 unit/contract tests, lint, and type checking pass. Signed-off-by: Kavya <58096384+maxmb1@users.noreply.github.com>
|
Warning This pull request is not mergeable via GitHub because a downstack PR is open. Once all requirements are satisfied, merge this PR as a stack on Graphite.
This stack of pull requests is managed by Graphite. Learn more about stacking. |
|
Closing this in favour of a new stack that lands the same tool with a smaller contract. Design: DistributedWorkers; tracking: REA-6343 ( Most of what is here carries over as is: the shared-memory block and What changes: the worker's methods become the model half's three names ( New PRs are linked from REA-6343 and REA-6345 as they open. Thanks for the groundwork; this stack stands on it. The new stack: #195 (transport and wire), #196 ( |
…#195) https://linear.app/reactor-team/issue/REA-6343 · Design: [DistributedWorkers](https://app.notion.com/p/reactor-technologies/DistributedWorkers-3decacd6a3968030b598d95d6df32ba0) First of three. Standalone: nothing else in the runtime imports `reactor_runtime.distributed`, and this stack does not depend on the `ReactorApp` stack (#187 to #192). Lifted from #134, #175, #176, #177, which it closes. ## Why The runner in #196 moves a model's input and result between processes. Both are the author's own objects and hold large arrays. Pickling them through a queue copies every byte; a fixed frame buffer forces a shape on the author. Neither is acceptable. ## What changed `pack()` pickles with protocol 5. Each flat buffer (a contiguous numpy array) goes into a shared-memory block; the pickle stays small and rides the queue as a header that names the block and the spans. `unpack()` attaches to whatever block the header names and copies the spans out. ```python slot = SharedSlot() # writer side, one per direction header = pack(result, slot) # ~1 KB on the queue; the arrays are in the block reader = SlotReader() # reader side result = unpack(header, reader) # the same object, arrays copied out ``` The block grows to the writer's high-water mark. Since the header names the block, growth needs no protocol. A torch tensor or a sliced array does not expose a flat buffer and pickles inline: slower, still correct, reported once per process. `SharedSlotAllocationFailed` names `--shm-size` and the pod's `/dev/shm` when a block cannot be created. Also here: the wire (`Generate`, `Reset`, `Shutdown`, `Loaded`, `Answer`) and the runner's errors (`WorkerCrashed`, `WorkerTimeout`, `RankDesync`). ## Verification `tests/unit/distributed/test_distributed_ipc.py`: a struct with a dict, str, enum, list, bytes, nested dataclass, `None`, and two arrays of different dtypes crosses intact with both arrays out of band; growth mid-call keeps what was already written and unlinks the old block; a non-contiguous array pickles inline and is reported once; the allocation error names the fix. `mise run lint`, `typecheck`, `test` green.

Why
Model authors should declare their worker and typed video output without reimplementing threading, session bookkeeping, and cleanup around WorkerGroup. This layer supplies that managed authoring path without introducing an engine-host dependency.
What Changed
Adds the experimental DistributedVideoModel adapter with worker configuration, control snapshots, frame-to-output mapping, pause/restart, and worker cleanup. Authoring and configuration tests cover the adapter contract, while the dependent example PR exercises the complete run-loop lifecycle. The scope remains single-node, single-session uint8 video, not arbitrary structured outputs or multi-session scheduling. Local Python 3.12 unit/contract tests, lint, and type checking pass.